Skip to main content

Bridging Streaming APIs Behind a REST Proxy

How to expose a simple JSON endpoint when your AI application does not speak REST JSON directly—for example SSE / text/event-stream, WebSockets, or gRPC—so you can integrate it with DynamoEval Custom Applications.

Interim workaround

DynamoEval Custom Applications currently integrate with model endpoints over REST and expect a single JSON request/response. Native support for SSE / event-stream, WebSockets, gRPC, and other non-REST protocols is not available yet.

Until that support lands, the recommended approach for anything other than REST JSON is to run a proxy outside of DynamoAI that:

  1. Accepts DynamoEval (or any REST client) over a normal JSON POST
  2. Calls your AI application over its native protocol (SSE, WebSocket, gRPC, and so on)
  3. Aggregates streamed or chunked output into one response when needed
  4. Returns a single JSON body that DynamoEval can consume

The problem​

Many LLM and chat APIs do not return one JSON blob. They may:

  • Stream tokens as Server-Sent Events (SSE) / text/event-stream
  • Speak WebSockets or gRPC instead of HTTP JSON POST
  • Otherwise expose a protocol DynamoEval cannot call directly

That works well for interactive UIs. It does not work well when:

  • Evaluation harnesses such as DynamoEval expect a single JSON response
  • Sync clients cannot consume event streams or binary RPC protocols
  • You need a stable contract even if the upstream protocol changes

Solution: put a REST proxy in front. Clients call a normal POST and get {"content": "..."}. The proxy talks to the upstream protocol, reads every chunk (or RPC response), concatenates the assistant text when streaming, and returns one JSON body.

The reference implementation below focuses on SSE, which is the most common case. The same adapter idea applies to WebSockets and gRPC: accept REST on the outside, speak the native protocol on the inside.

Client ──POST /proxy──► Proxy ──stream POST──► Upstream (SSE)
│
▼
consume stream
aggregate deltas
│
Client ◄──{"content": "..."}──┘

After the proxy is running, register its REST URL as the remote model endpoint in DynamoEval using the standard Custom API Language Models flow (including JSONata transformations if needed). See Custom API Language Model Requirements and Custom Systems Overview for prerequisites, authentication, and transformation basics.

Architecture​

LayerRole
Public REST endpointAccepts JSON, returns aggregated JSON
Stream forwarderOpens an HTTP stream to the upstream with httpx
Stream consumerParses each SSE event and extracts text deltas (application-specific)
Error mapperMaps timeouts and upstream HTTP errors to clear status codes

Reference stack used in this guide: FastAPI + httpx (async streaming). You can implement the same pattern in any language that supports HTTP streaming.

Authentication​

How you handle auth is up to you. Two common patterns:

  1. Auth at the proxy — Configure DynamoEval’s Custom AI system with credentials for the proxy. The proxy validates those credentials, then authenticates itself to the upstream AI application (for example with a service token it holds).
  2. Pass-through auth — Configure DynamoEval’s Custom AI system with the AI application’s real credentials. The proxy forwards those headers unchanged to the upstream endpoint.

Either approach works with DynamoEval’s supported auth modes (No Auth, Bearer Token, API Key). Pick the model that matches your security and ops preferences.

Step 1 — Public REST endpoint​

Expose a normal POST route. Callers see REST; they never see the stream.

from fastapi import FastAPI, Request, HTTPException

app = FastAPI(title="Streaming Bridge Proxy")

UPSTREAM_URL = "https://your-streaming-api.example.com/v1/chat"

@app.post("/proxy")
async def proxy_chat(request: Request):
try:
return await make_sse_request(UPSTREAM_URL, request)
except HTTPException:
raise
except Exception as e:
return {"content": f"UNCAUGHT_ERROR: {type(e).__name__}: {e}"}

Response shape:

{
"content": "The full assistant reply, assembled from every stream delta."
}

Map this shape to DynamoEval’s expected response format with a JSONata response transformation if needed—for example [content].

Step 2 — Forward the request as a stream​

Use httpx.AsyncClient.stream() so the connection stays open while chunks arrive. Do not use a one-shot client.post() that buffers the whole body before you can parse it.

import httpx
from fastapi import Request, HTTPException, Response

TIMEOUT_IN_S = 180.0

async def make_sse_request(url: str, request: Request):
try:
client_body = await request.json()
except Exception:
raise HTTPException(status_code=400, detail="Invalid JSON body received.")

async with httpx.AsyncClient() as client:
try:
async with client.stream(
"POST",
url,
json=client_body,
headers=dict(request.headers),
timeout=TIMEOUT_IN_S,
) as response:
if response.status_code >= 400:
error_body = (await response.aread()).decode("utf-8", errors="replace")
raise HTTPException(status_code=response.status_code, detail=error_body)

content = await consume_stream(response)
return {"content": content}

except httpx.ConnectTimeout:
return Response(status_code=502, content="LLM Connection Timeout")
except httpx.ReadTimeout:
return Response(status_code=504, content="LLM Read Timeout")
except httpx.ConnectError:
return Response(status_code=502, content="LLM Connection Error")

Prefer a long timeout (for example, 180s). Streams finish when generation finishes, not when the first byte arrives.

Step 3 — Consume and aggregate the stream​

Application-specific parsing

The consume_stream logic below is an example for one AI application’s SSE event shape. How you parse and aggregate the stream depends entirely on how your AI application streams responses—event names, JSON fields, and framing (data: lines, keep-alives, and so on) will differ. Treat the functions here as a starting point and adapt the extractor to your upstream schema.

SSE typically sends lines like data: {...}. Strip the prefix, parse JSON, and pull out text deltas.

Example event payload (illustrative only):

{"obj": {"type": "message_delta", "content": " Hi"}}

Consumer (example—rewrite extract_delta_content for your API):

import json
import logging

logger = logging.getLogger(__name__)

def extract_delta_content(payload: dict) -> str | None:
"""Pull assistant text from one SSE event. Adapt to your upstream schema."""
obj = payload.get("obj")
if not isinstance(obj, dict):
return None
if obj.get("type") != "message_delta":
return None
content = obj.get("content")
return content if isinstance(content, str) else None


async def consume_stream(response: httpx.Response) -> str:
"""Read SSE lines and concatenate text deltas for this application."""
content_parts: list[str] = []

async for line in response.aiter_lines():
if not line or line.startswith(":"):
continue # keep-alives / comments
if line.startswith("data:"):
line = line[5:].strip()
if not line.strip():
continue
try:
payload = json.loads(line)
except json.JSONDecodeError:
logger.warning("Skipping non-JSON stream line: %r", line)
continue

chunk = extract_delta_content(payload)
if chunk:
content_parts.append(chunk)

return "".join(content_parts)

The aggregation idea is always the same: parse → extract text → join → return once. Only the parsing and field extraction change per application.

Step 4 — Error handling map​

FailureTypical statusMeaning
Invalid client JSON400Bad request body
Upstream 4xx / 5xxSame as upstreamPass through detail when safe
Connect timeout502Couldn’t reach the model service
Read timeout504Stream didn’t finish in time
Connect error502Network / DNS / TLS failure

Check upstream status before consuming the body so error payloads are not treated as chat deltas.

End-to-end flow​

1. Client POSTs JSON to /proxy
2. Proxy validates JSON
3. Proxy opens streamed POST to upstream (headers passed through)
4. If status >= 400 → raise with upstream body
5. Else read SSE events until the stream ends
6. For each event, extract text deltas using your app-specific parser
7. Return {"content": "<joined text>"}

Why this pattern works well​

For callers

  • One request, one JSON response
  • Works with DynamoEval, scripts, and sync HTTP clients
  • Stable contract even if upstream streaming format evolves

For operators

  • Flexible auth (proxy-owned credentials or pass-through to the AI application)
  • Timeouts and connection errors map to clear HTTP codes
  • Logging can sit at the proxy without changing clients

Trade-off

  • Clients wait until the full answer is generated (no token-by-token UI). That is intentional when the goal is a complete REST response for evaluation, not live streaming.

Minimal dependencies​

fastapi>=0.104.0
uvicorn[standard]>=0.24.0
httpx>=0.25.0

Adapting this to your API​

  1. Point UPSTREAM_URL at your streaming chat endpoint.
  2. Rewrite extract_delta_content() / consume_stream() to match your SSE event schema (delta, choices[0].delta.content, custom event types, and so on).
  3. Choose an auth pattern (proxy auth vs pass-through) and configure DynamoEval accordingly.
  4. Keep aggregation and error mapping; the stream consumer is what usually changes.
  5. Point DynamoEval’s custom model remote_model_endpoint at the proxy’s REST URL and configure auth / JSONata as usual.

Summary​

A streaming-to-REST proxy is an adapter:

  1. Accept a normal JSON POST
  2. Open an HTTP stream upstream
  3. Parse and concatenate text deltas (using logic specific to your AI application)
  4. Return one aggregated JSON body

With FastAPI and httpx’s client.stream(), that bridge is a few focused functions—and it lets DynamoEval (and other tools that only speak REST) talk to backends that only speak SSE.